Java 响应式编程 - Reactive Streams 规范
在现代高并发、高吞吐量的互联网应用中,传统的一个请求对应一个线程模型很容易在 I/O 阻塞(如数据库查询、远程调用)时导致线程资源耗尽。为了解决这个问题,响应式编程(Reactive Programming)应运而生。而 Reactive Streams(响应式流)规范,则是整个 Java 生态响应式编程的基石。
什么是 Reactive Streams 规范
Reactive Streams 是一套面向流式数据处理、具有 非阻塞背压(Non-blocking Backpressure)的异步组件标准。简单来说,它定义了一套 API 规范,让不同的框架(如 Spring WebFlux、RxJava、Akka Streams)可以互相兼容,共同解决 “上游生产数据太快,下游消费太慢,导致内存溢出(OOM)” 的核心痛点。
Reactive Streams 的核心特征就是两个字——背压。在传统的流处理中,如果上游(Publisher)发送数据的速度远快于下游(Subscriber)处理的速度,下游的内存缓冲区就会被撑爆。背压机制允许下游订阅者显式地告诉上游:“我现在只能处理 3 个元素,请不要发多了。” 从而实现了流量控制,保证了系统的稳定性。
Reactive Streams 官网 reactive-streams.org。它的 Java 实现非常精简,在 Java 9 之中被直接引入到了 java.util.concurrent.Flow 类中。
该规范极其精简,核心只有 4 个接口(四大核心组件):
1 | // 1. 发布者:数据的源头,负责生产数据 |
Reactive Streams 的原理和流程
并发模型
Reactive Streams 规范本身没有绑定任何特定的线程模型,它是一套协议。我们这里说的其原理和工作流程对应到它的落地到具体实现,比如 Project Reactor 和 Spring WebFlux 上。
第一,在多线程并发下,下游可能在主线程请求数据,上游可能在定时器线程产生数据。为了保证高并发下的绝对安全且不阻塞,Subscription 内部通常使用 CAS 乐观锁和原子变量(如 AtomicLong、AtomicReference)来维护两个核心状态:
- requested(请求计数器):下游请求的剩余安全配额。每次下游 request(n),利用 CAS 执行 requested += n;每次上游 onNext(),执行 requested–。当 requested == 0 时,上游必须彻底停止发送,原地待命。
- cancelled(取消标志):布尔原子量。一旦变为 true,后续所有 onNext 直接被丢弃并清理内存。
第二,在传统多线程编程中,我们需要自己处理 ExecutorService、CountDownLatch、线程间共享队列的可见性、死锁等问题。Reactive Streams 框架(以 Project Reactor 为例)通过将 “线程切换” 抽象为声明式的操作符,彻底屏蔽了这些复杂性:
- publishOn(Scheduler):改变下游执行的线程池。
- subscribeOn(Scheduler):改变上游数据源加载时执行的线程池。
第三,当你写下 flux.publishOn(Schedulers.boundedElastic()).map(…) 时,Reactor 底层实际上利用了异步环形队列(Ring Buffer / Queue)和事件驱动回调:
- 无锁队列隔离:publishOn 内部会偷偷塞入一个 Processor 或一个高性能的单生产者-单消费者(SPSC)异步无锁队列。
- 线程切换点:上游线程(假设是 A 线程)生产了数据,它不直接调用下游的 onNext,而是把数据塞进这个队列,并向目标线程池(BoundedElastic 线程池中的 B 线程)提交一个执行任务。
- 消费者线程接管:B 线程苏醒后,从队列中源源不断地取出数据,然后高效率地执行下游的 map 和 subscribe 逻辑。
这一切都被封装在操作符内部,开发者看到的只是优雅的链式调用。在应用运行时,底层的线程开多少、怎么开,完全取决于宿主服务器(如 Netty)和调度器(Schedulers)的配置。我们以 Spring WebFlux + Netty + Reactor 这一黄金组合来拆解:
① 核心网络 I/O 线程池:Netty EventLoopGroup:当你启动一个 WebFlux 应用时,底层默认运行着一个 Netty 服务器。
- 开了多少线程?:默认是 CPU 核心数 x 2(通常被称为 reactor-http-nio-X 线程)。
- 怎么工作?:这些是极度宝贵的非阻塞 Worker 线程。它们利用操作系统的多路复用(如 Linux 的 epoll),一个线程就能同时管理成千上万个网络连接的 Socket 读写。
- 致命禁忌:绝对不能在这些线程里执行任何阻塞操作(如 JDBC、Thread.sleep())。一旦阻塞,Netty 的事件循环就会卡死,整个服务器的吞吐量会瞬间雪崩。
② Reactor 的内置调度器(线程池类型)
为了处理不同类型的业务,Reactor 提供了几类开箱即用的线程池(Schedulers):

核心工作流程
1 | Publisher Subscriber |
- 建立连接 (subscribe):Subscriber 调用 Publisher.subscribe(Subscriber)。
- 初始化并传递令牌 (onSubscribe):Publisher 创建一个 Subscription 对象(常命名为 ScalarSubscription 或 FluxHide.SuppressSubscription),通过 Subscriber.onSubscribe(Subscription) 传给下游。此时,上游不会发送任何数据。
- 下游按需索要 (request):下游 Subscriber 根据自身处理能力,调用 Subscription.request(n)。这一步通常是异步的。
- 上游生产并推送 (onNext):上游收到 request(n) 信号后,开始生产数据,并调用 Subscriber.onNext(T) 传递数据。上游发送的数据量绝对不能超过 n。
- 循环或结束 (onComplete / onError):下游消费完这 $n$ 个数据后,如果还需要,就继续调用 request(m);如果数据发完了,上游调用 onComplete() 关闭流。
一个完整的请求
我们再来看一个 WebFlux 收到请求并调用 R2DBC(Reactive Relational Database Connectivity)的底层真实全链路线程流转轨迹:
1 | [用户浏览器] |
你可能会有疑惑:整个过程中,查询数据库这一步不是阻塞的吗?如果是异步调用,数据库返回数据又是怎么返回 前端的呢?绝大多数从传统 Servlet(Spring MVC)转到响应式编程的开发者,都会在这里卡壳:数据库查询怎么可能不阻塞?哪怕用非阻塞,数据回来后又是怎么神不知鬼不觉地塞回前端的?这里隐藏着两个底层的关键技术:R2DBC(驱动层的革新) 和 Netty Channel Pipeline(网络层的复用)。
传统的 JDBC(如 mysql-connector-java)之所以是阻塞的,是因为它的底层通信使用了传统的 BIO(Blocking I/O)。当 Java 线程执行 resultSet.next() 时,线程会发出一个 TCP 请求,然后原地死等数据库服务器的响应。此时操作系统会将该线程挂起(Block),直至收到数据。
而响应式编程使用的 R2DBC,其底层彻底抛弃了 BIO,改用 NIO(Non-blocking I/O,通常基于 Netty)。底层真实的执行过程是:
- 发送请求后立即撒手:当你的业务代码调用数据库查询时,Java 线程(Netty Worker 线程)将 SQL 语句打包成一个 TCP 数据包,通过操作系统的非阻塞 Socket 发送出去。
- 不等待,直接返回:发送完毕后,该线程绝对不会在原地等待数据库返回结果。它在驱动层注册一个 “数据准备就绪” 的回调事件(Event Listener),然后这个线程就立刻结束当前任务,去高高兴兴地接待下一个用户的 HTTP 请求了。
- 内核接管:此时,没有任何 Java 线程在为这个数据库查询买单。只有操作系统的内核(如 Linux 的 epoll)在盯着网卡。
既然发送请求的线程已经跑去干别的事了,那么当数据库执行完 SQL,把数据顺着网络线缆丢回来的时候,数据是怎么一步步回到前端浏览器的呢?这里全靠操作系统的事件通知机制和 Reactive Streams 的链式订阅关系来完成接力。整个过程就像一个完美的多米诺骨牌效应:
- 骨牌 1:网卡收到数据,内核唤醒 Netty。数据库把数据计算好后,通过 TCP 回传。服务器的网卡收到数据包,引发一个硬件中断。Linux 内核发现这个 Socket 有数据进来了,于是通知正在执行 epoll_wait 的 Netty 线程池(可能是刚才那个线程,也可能是同池子里的另一个空闲 Worker 线程,我们假设是 reactor-http-nio-2)。
- 骨牌 2:驱动程序触发 Reactive Streams 信号。reactor-http-nio-2 线程领到任务,开始读取网卡缓冲区里的数据库原始字节流。R2DBC 驱动将这些字节解析成 Java 对象。此时,驱动程序作为 Publisher(发布者),开始向下游触发 Reactive Streams 标准信号:
- 它调用下游的 Subscription.request(n) 看需不需要数据。
- 确认需要后,调用 Subscriber.onNext(data) 将转换好的数据行(Row)推送出去。
- 骨牌 3:业务操作符链条被 “串路激活”。数据一旦进入 onNext,你之前写的那些声明式代码(如 .map(user -> …), .filter(…))就开始在 reactor-http-nio-2 线程上流水线式地执行。这些代码并不是在请求时执行的,而是在数据回来的这一刻,被数据流“冲刷”激活的。
- 骨牌 4:Netty 写入网络缓冲区,直达前端。当数据流一路走到最下游(WebFlux 框架层)时,WebFlux 的底层订阅者会把这些数据对象序列化为 JSON 字符串,然后直接调用 Netty 的 channel.writeAndFlush()。 Netty 再次利用非阻塞 I/O,把 JSON 字节写入到与前端浏览器建立的那个 Socket 管道(Channel)中。网卡将数据发送出去,前端浏览器收到响应,渲染页面。
通过这种 “事件驱动、收到通知才干活” 的方式,哪怕数据库查询需要耗时 3 秒,Java 服务器在这 3 秒内也不需要耗费任何一个线程去死等,从而释放了恐怖的并发吞吐能力。
两个具体的测试案例
基础案例
以下是一个只包含发布者和订阅者的最简单的案例:
1 | import java.io.IOException; |
包含 Processor 的案例
Processor(处理者)同时继承了 Subscriber 和 Publisher,这意味着它扮演了一个“双面间谍” 的角色:
- 对于上游,它是一个订阅者,负责接收原始数据;
- 对于下游,它是一个发布者,负责把加工、过滤、转换后的数据推送出去。
下面的案例中:Processor 拦截上游的原始订单,自动过滤掉 100 元以下的“小额订单”,并将剩余的高端订单包装成 “VIP专属订单” 再投递给下游的风控系统。这个案例完整展示了底层的链式接力机制:
Publisher (数据源) —–> Processor (加工过滤) —–> Subscriber (最终消费)
1 | import lombok.AllArgsConstructor; |
1 | [Subscriber] 成功订阅,请求 1 个高端订单... |
关于响应式编程的学习
响应式编程(Reactive Programming)也许你刚开始觉得很晦涩,是因为它彻底颠覆了我们写了十几年的命令式编程(Imperative Programming)思维。过去我们习惯了先 A 再 B 卡住等结果,而响应式关注的是数据怎么流动,怎么对变化做出响应。为了让你彻底看清,我们跳出复杂的 API,用一条 “由浅入深、层层递进” 的黄金主线来帮你把零散的知识点串联起来。
- 第一站:思维重构(从拉模型到推模型)
- 传统思维(拉模型 - Pull):你作为老板(主线程),让员工 A 去数据库拿数据。A 去了,你就在办公室坐着死等(阻塞),A 回来之前你什么都干不了。
- 响应式思维(推模型 - Push):你(主线程)不再等 A,而是构建一条生产流水线/自来水管。你跟 A 说:“你拿到数据后,顺着这条管道丢进来。” 交代完之后,你立刻去干别的事。当数据真正顺着管道流下来的时候,流水线上的各个节点自动触发加工。
- 响应式编程 = 声明式数据流(Data Stream) + 变化传递(Propagation of change)。你在写代码时,不是在下达执行命令,而是在“搭管道”。
- 第二站:统一标准(Reactive Streams 规范)
- 当你理解了管道思维,就会发现一个问题:上游放水太快,下游池子太小,水会溢出来(也就是内存溢出 OOM)。于是,Java 制定了 Reactive Streams 规范。
- 规范极其简单,它用 4 个接口规定了管道两端和中间的玩法:
- Publisher(发布者 / 水龙头):负责生产数据。
- Subscriber(订阅者 / 盛水的水桶):负责消费数据。
- Subscription(订阅关系 / 水阀开关):最核心的灵魂! 下游调用 request(n) 告诉上游:“我现在有空,给我放 $n$ 滴水。” 上游放完 $n$ 滴立刻停下,等下次通知。这就是背压机制。
- Processor(处理器 / 过滤器):既是水桶又是水龙头,放在中间净水。
- 规范不包含任何复杂的线程池或高深技术,它只是通过 request(n) 这一小步,完美解决了 “上游太快饿死/撑死下游” 的千古难题。
- 第三站:乐高组装(Project Reactor 核心)
- 规范只是设计图纸,无法直接在生产中使用。Pivotal 团队(Spring 亲爹)根据这个规范写出了工业级实现——Project Reactor。它把晦涩的接口变成了如同乐高积木般好玩的工具。它提供了两个核心管道工:
- Mono\<T>:只会流出 0 或 1 个元素的特制管道。
- Flux\<T>:可以流出 0 到 N 个无穷元素的标准管道。
- 它最大的魅力在于提供了几百个操作符(Operators),就像各式各样的管道配件:
- filter():过滤掉不合格的水(数据)。
- map():把清水变成纯净水(数据类型转换)
- publishOn() / subscribeOn():魔法切换器。水流到这里,自动切换到另一个动力泵(线程池)继续推,开发者完全不用操心底层的 Thread 和 Lock。
- 完全不要去死记几百个操作符!记住它们只是管道配件,随用随查。你只需要用链式调用把 Flux.just(1, 2).map(…).filter(…) 串起来即可。
- 规范只是设计图纸,无法直接在生产中使用。Pivotal 团队(Spring 亲爹)根据这个规范写出了工业级实现——Project Reactor。它把晦涩的接口变成了如同乐高积木般好玩的工具。它提供了两个核心管道工:
- 第四站:实战落地(Spring WebFlux)
- 当你学会了组装管道,最后一步就是用它来重构你的 Web 应用。这就是 Spring WebFlux。它彻底抛弃了 Tomcat 那种“一个请求死死绑定一个线程”的模式,改用 Netty(基于事件驱动的非阻塞服务器)。
- 传统 Tomcat:1000 个并发请求,必须要 1000 个物理线程。如果数据库很慢,1000 个线程全部卡死。
- WebFlux + Netty:可能一共就 4 个核心线程。遇到 I/O 阻塞(比如查数据库),这 4 个线程“把请求往操作系统的异步队列里一丢”,立刻转头去接第 5 个请求。等到数据库数据回来了,操作系统发个通知,4 个线程中的某一个再回来把数据顺着管道推给前端。
- WebFlux 的并发能力不是靠增加线程实现的,而是靠 “不让线程闲着” 实现的。
- 当你学会了组装管道,最后一步就是用它来重构你的 Web 应用。这就是 Spring WebFlux。它彻底抛弃了 Tomcat 那种“一个请求死死绑定一个线程”的模式,改用 Netty(基于事件驱动的非阻塞服务器)。
- 避坑建议:
- 忘记 ThreadLocal:不要再奢望用 ThreadLocal 传递登录信息或日志 MDC 令牌了,因为响应式下,你的代码可能前 1 秒在线程 A 执行,后 1 秒就在线程 B 执行了。要习惯使用 Reactor 的 Context。
- 绝对不要在管道里“下毒”:在 Flux/Mono 的链条里,绝对绝对不能调用阻塞方法(比如老版 HttpClient.execute()、Thread.sleep() 或传统的 JDBC)。只要有一个阻塞,整个响应式线程池就会瘫痪。
- 记住“不订阅,什么都不会发生”:比如 Mono\<User> user = r2dbcRepository.findById(1L) ,这行代码执行了,但由于最后没有 subscribe(),数据库连查都不会查,它只是一张设计图!
按照 “重构思维 $\rightarrow$ 搞懂背压接口 $\rightarrow$ 熟悉 Flux/Mono 链式操作 $\rightarrow$ 结合 WebFlux” 这条主线去学,你会发现那些晦涩的 API 背后,其实是一套极其优雅、有条不紊的现代化工业流水线。